FlowStreamType

FlowStream

FlowStream<'env, 'error, 'value> describes a cold, pull-based sequence of effectful values. It uses the same environment, typed failure channel, cancellation, and runtime scopes as Flow.

Choose FlowStream when values should be handled incrementally instead of collected before work begins. Typical sources include process output, paginated APIs, subscriptions, and host-specific file or network adapters.

A stream does nothing until a terminal operation such as runFold, runForEachFlow, or runCollect produces a Flow and that Flow is run. Pulling provides backpressure: downstream requests each next value, so upstream cannot run ahead without an operator explicitly introducing bounded concurrency. The time-based operators (groupedWithin, throttle, debounce, switchMapFlow) and buffer are such operators: they read ahead into a queue so they can react to time or to a newer value.

A stream that reads a file, socket, or cursor should own it with using, which releases the resource however consumption ends. To take a single result without draining the stream, use runTryHead.

Streams connect to the concurrency types in both directions. fromDequeue reads a queue or hub subscription until it is shut down and drained, and fromHub subscribes for the life of the stream. runIntoQueue and runIntoHub feed a stream into them. mergePar runs several streams concurrently, and fromSchedule turns a schedule into a tick stream. See the Queue and Hub guides.

Start here

The Streams guide introduces the model. Getting started builds a complete pipeline, while Batching and parallelism explains the scheduling difference between strict batches and mapFlowPar.

Use incremental terminal operations for large or unbounded streams. runCollect deliberately retains every emitted value and should be reserved for known finite inputs.

Specification

Kind
Type
Members
47
Examples
0
localEnvSignature
FlowStream.localEnv mapping (FlowStream operation)
unfoldFlowSignature
FlowStream.unfoldFlow step initialState
fromDequeueSignature
FlowStream.fromDequeue queue
fromHubSignature
FlowStream.fromHub strategy hub
fromScheduleSignature
FlowStream.fromSchedule schedule

Summary

NameSignatureSynopsis
Union cases
FlowStreamFlowStream 'env -> CancellationToken -> Execution, 'error>No description available.
Operations
localEnvFlowStream.localEnv mapping (FlowStream operation)Runs a stream with an environment derived from its caller's environment.
unfoldFlowFlowStream.unfoldFlow step initialStateCreates a cold stream by repeatedly running an effectful state transition.
fromDequeueFlowStream.fromDequeue queueCreates a stream that takes values from a queue or hub subscription until it is shut down and drained.
fromHubFlowStream.fromHub strategy hubCreates a stream of the values published to a hub, subscribing for the life of the stream.
fromScheduleFlowStream.fromSchedule scheduleCreates a stream that emits a schedule's outputs, each after the delay the schedule chooses.
runIntoQueueFlowStream.runIntoQueue queue streamRuns a stream and offers every value to a queue, waiting while a bounded queue is full.
runIntoHubFlowStream.runIntoHub hub streamRuns a stream and publishes every value to a hub.
bufferFlowStream.buffer strategy streamRuns the upstream ahead of its consumer in its own fiber, buffering values with strategy.
mergeParFlowStream.mergePar streamsRuns several streams concurrently and emits their values as they arrive.
fromSeqFlowStream.fromSeq valuesCreates a stream from a synchronous sequence of values.
emptyFlowStream.empty Creates an empty stream.
singletonFlowStream.singleton valueCreates a stream containing one value.
fromFlowFlowStream.fromFlow flowCreates a one-element stream from an effectful value.
runForEachFlowStream.runForEach action (FlowStream op)Executes the stream and performs a synchronous action for each successful value.
mapFlowStream.map f streamTransforms the successful values of a stream using the provided function.
mapErrorFlowStream.mapError mapper streamTransforms the typed error channel of a stream.
filterFlowStream.filter predicate streamKeeps values that satisfy a predicate.
chooseFlowStream.choose chooser streamMaps and filters values in one operation.
tapFlowFlowStream.tapFlow action streamRuns an effect for each value before emitting the original value.
mapFlowFlowStream.mapFlow mapper streamTransforms every value with a Flow effect.
mapFlowParFlowStream.mapFlowPar parallelism mapper streamMaps values with a continuously replenished, bounded set of child fibers.
mapFlowParUsingFlowStream.mapFlowParUsing parallelism resource mapper streamLike mapFlowPar, giving each running mapping a resource from a pool that holds at most parallelism of them.
groupedWithinFlowStream.groupedWithin size window streamGroups values into lists of at most size, emitting early when window passes.
debounceFlowStream.debounce quiet streamEmits a value only once quiet passes without a newer one.
throttleFlowStream.throttle interval streamEmits at most one value per interval, keeping the latest value when faster.
switchMapFlowFlowStream.switchMapFlow mapper streamMaps each value to a flow, interrupting the previous flow when a newer value arrives.
chunkBySizeFlowStream.chunkBySize size streamGroups consecutive values into non-empty lists of at most size elements.
takeFlowStream.take count streamEmits at most count values.
skipFlowStream.skip count streamSkips the first count values.
takeWhileFlowStream.takeWhile predicate streamEmits values while a predicate remains true.
skipWhileFlowStream.skipWhile predicate streamSkips values while a predicate remains true.
indexedFlowStream.indexed streamEmits each value paired with its zero-based index.
scanFlowStream.scan folder initial streamEmits successive accumulator states.
distinctUntilChangedByFlowStream.distinctUntilChangedBy projection streamSuppresses consecutive duplicate values according to a projection.
appendFlowStream.append (FlowStream right) (FlowStream left)Concatenates two streams, evaluating the second only after the first ends.
collectFlowStream.collect mapper streamMaps each value to a stream and concatenates the resulting streams.
zipFlowStream.zip rightStream leftStreamPairs values from two streams until either stream ends.
runFoldFlowStream.runFold folder initial streamFolds a stream into one value inside Flow.
runCollectFlowStream.runCollect streamCollects all emitted values into a list.
runDrainFlowStream.runDrain streamConsumes a stream and ignores its values.
runForEachFlowFlowStream.runForEachFlow action streamRuns an effectful action for every stream value.
repeatFlowFlowStream.repeatFlow flowCreates a stream that emits the result of running flow again for every pull, forever.
usingFlowStream.using resource buildCreates a stream that owns a resource: acquired when consumption starts, released when it ends.
runTryHeadFlowStream.runTryHead streamReturns the first value, or None for an empty stream, then stops the stream.
runTryLastFlowStream.runTryLast streamConsumes the whole stream and returns its last value, or None for an empty stream.
runCountFlowStream.runCount streamConsumes the whole stream and returns how many values it emitted.

Union cases

kind:member

FlowStream

FlowStream 'env -> CancellationToken -> Execution<StreamStep<'value, 'error>, 'error>
Member

Parameters

NameTypeDescription
Item'env -> CancellationToken -> Execution<StreamStep<'value, 'error>, 'error>

Returns

unit

Operations

kind:member

localEnv

FlowStream.localEnv mapping (FlowStream operation)
Member
Runs a stream with an environment derived from its caller's environment.

Parameters

NameTypeDescription
mapping'outerEnvironment -> 'innerEnvironment
FlowStream operationFlowStream<'innerEnvironment, 'error, 'value>

Returns

FlowStream<'outerEnvironment, 'error, 'value>
kind:member

unfoldFlow

FlowStream.unfoldFlow step initialState
Member
Creates a cold stream by repeatedly running an effectful state transition.

Parameters

NameTypeDescription
step'state -> Flow<'env, 'error, ('value * 'state) option>Returns Some(value, nextState) or None when the stream is complete.
initialState'stateThe state used for the first pull.

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

FlowStream.unfoldFlow (fun n -> Flow.ok (if n < 3 then Some(n, n + 1) else None)) 0
kind:member

fromDequeue

FlowStream.fromDequeue queue
Member
Creates a stream that takes values from a queue or hub subscription until it is shut down and drained.

Parameters

NameTypeDescription
queueDequeue<'value>

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

flow {
    let! (jobs: Queue<string>) = Queue.bounded 8
    do! jobs |> Queue.offerAll [ "a"; "b" ] |> Flow.ignore
    do! Dequeue.shutdown jobs
    return! jobs |> FlowStream.fromDequeue |> FlowStream.runCollect
}
kind:member

fromHub

FlowStream.fromHub strategy hub
Member
Creates a stream of the values published to a hub, subscribing for the life of the stream.

Parameters

NameTypeDescription
strategyQueueStrategy
hubHub<'value>

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

readings
|> FlowStream.fromHub (QueueStrategy.Sliding 1)
|> FlowStream.runForEach display
kind:member

fromSchedule

FlowStream.fromSchedule schedule
Member
Creates a stream that emits a schedule's outputs, each after the delay the schedule chooses.

Parameters

NameTypeDescription
scheduleSchedule<'env, unit, 'output>

Returns

FlowStream<'env, 'error, 'output>

Verification Examples

// A tick every 50 ms, aligned to when the stream started.
Schedule.fixedRate (TimeSpan.FromMilliseconds 50.0) |> FlowStream.fromSchedule
kind:member

runIntoQueue

FlowStream.runIntoQueue queue stream
Member
Runs a stream and offers every value to a queue, waiting while a bounded queue is full.

Parameters

NameTypeDescription
queueQueue<'value>
streamFlowStream<'env, 'error, 'value>

Returns

Flow<'env, 'error, unit>

Verification Examples

deviceReadings |> FlowStream.runIntoQueue controlInputs
kind:member

runIntoHub

FlowStream.runIntoHub hub stream
Member
Runs a stream and publishes every value to a hub.

Parameters

NameTypeDescription
hubHub<'value>
streamFlowStream<'env, 'error, 'value>

Returns

Flow<'env, 'error, unit>

Verification Examples

sensor |> FlowStream.runIntoHub readings
kind:member

buffer

FlowStream.buffer strategy stream
Member
Runs the upstream ahead of its consumer in its own fiber, buffering values with strategy.

Parameters

NameTypeDescription
strategyQueueStrategy
streamFlowStream<'env, 'error, 'value>

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

samples |> FlowStream.buffer (QueueStrategy.Sliding 1) |> FlowStream.runForEachFlow render
kind:member

mergePar

FlowStream.mergePar streams
Member
Runs several streams concurrently and emits their values as they arrive.

Parameters

NameTypeDescription
streamsFlowStream<'env, 'error, 'value> list

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

[ readerA; readerB; readerC ] |> FlowStream.mergePar |> FlowStream.runIntoQueue controlInputs
kind:member

fromSeq

FlowStream.fromSeq values
Member
Creates a stream from a synchronous sequence of values.

Parameters

NameTypeDescription
values'value seqThe sequence of values to be emitted by the stream.

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

FlowStream.fromSeq [1..10]
|> FlowStream.runCollect
|> Flow.run ()
kind:member

empty

FlowStream.empty
Member
Creates an empty stream.

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

FlowStream.empty<unit, string, int>
kind:member

singleton

FlowStream.singleton value
Member
Creates a stream containing one value.

Parameters

NameTypeDescription
value'value

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

FlowStream.singleton 42
kind:member

fromFlow

FlowStream.fromFlow flow
Member
Creates a one-element stream from an effectful value.

Parameters

NameTypeDescription
flowFlow<'env, 'error, 'value>

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

FlowStream.fromFlow (Flow.ok 42)
kind:member

runForEach

FlowStream.runForEach action (FlowStream op)
Member
Executes the stream and performs a synchronous action for each successful value.

Parameters

NameTypeDescription
action'value -> unitThe function to execute for each value emitted by the stream.
FlowStream opFlowStream<'env, 'error, 'value>

Returns

Flow<'env, 'error, unit>

Verification Examples

FlowStream.fromSeq ["a"; "b"; "c"]
|> FlowStream.runForEach (printfn "%s")
|> Flow.run ()
kind:member

map

FlowStream.map f stream
Member
Transforms the successful values of a stream using the provided function.

Parameters

NameTypeDescription
f'v -> 'wThe function to transform each value.
streamFlowStream<'env, 'error, 'v>The stream whose values should be transformed.

Returns

FlowStream<'env, 'error, 'w>

Verification Examples

let stream = FlowStream.fromSeq [1; 2; 3] |> FlowStream.map (fun n -> n * 2)
kind:member

mapError

FlowStream.mapError mapper stream
Member
Transforms the typed error channel of a stream.

Parameters

NameTypeDescription
mapper'error -> 'nextError
streamFlowStream<'env, 'error, 'value>

Returns

FlowStream<'env, 'nextError, 'value>

Verification Examples

stream |> FlowStream.mapError DomainError
kind:member

filter

FlowStream.filter predicate stream
Member
Keeps values that satisfy a predicate.

Parameters

NameTypeDescription
predicate'a -> bool
streamFlowStream<'b, 'c, 'a>

Returns

FlowStream<'b, 'c, 'a>

Verification Examples

stream |> FlowStream.filter (fun value -> value > 0)
kind:member

choose

FlowStream.choose chooser stream
Member
Maps and filters values in one operation.

Parameters

NameTypeDescription
chooser'a -> 'b option
streamFlowStream<'c, 'd, 'a>

Returns

FlowStream<'c, 'd, 'b>

Verification Examples

stream |> FlowStream.choose id
kind:member

tapFlow

FlowStream.tapFlow action stream
Member
Runs an effect for each value before emitting the original value.

Parameters

NameTypeDescription
action'a -> Flow<'b, 'c, unit>
streamFlowStream<'b, 'c, 'a>

Returns

FlowStream<'b, 'c, 'a>

Verification Examples

stream |> FlowStream.tapFlow logValue
kind:member

mapFlow

FlowStream.mapFlow mapper stream
Member
Transforms every value with a Flow effect.

Parameters

NameTypeDescription
mapper'a -> Flow<'b, 'c, 'd>
streamFlowStream<'b, 'c, 'a>

Returns

FlowStream<'b, 'c, 'd>

Verification Examples

ids |> FlowStream.mapFlow load
kind:member

mapFlowPar

FlowStream.mapFlowPar parallelism mapper stream
Member
Maps values with a continuously replenished, bounded set of child fibers.

Parameters

NameTypeDescription
parallelismParallelism
mapper'a -> Flow<'b, 'c, 'd>
streamFlowStream<'b, 'c, 'a>

Returns

FlowStream<'b, 'c, 'd>
kind:member

mapFlowParUsing

FlowStream.mapFlowParUsing parallelism resource mapper stream
Member
Like mapFlowPar, giving each running mapping a resource from a pool that holds at most parallelism of them.

Parameters

NameTypeDescription
parallelismParallelism
resourceResource<'env, 'error, 'resource>
mapper'resource -> 'value -> Flow<'env, 'error, 'next>
streamFlowStream<'env, 'error, 'value>

Returns

FlowStream<'env, 'error, 'next>

Verification Examples

commits |> FlowStream.mapFlowParUsing (Parallelism.bounded 4) openReader (fun reader commit -> diff reader commit)
kind:member

groupedWithin

FlowStream.groupedWithin size window stream
Member
Groups values into lists of at most size, emitting early when window passes.

Parameters

NameTypeDescription
sizeint
windowTimeSpan
streamFlowStream<'env, 'error, 'value>

Returns

FlowStream<'env, 'error, 'value list>

Verification Examples

events |> FlowStream.groupedWithin 100 (TimeSpan.FromSeconds 1.0)
kind:member

debounce

FlowStream.debounce quiet stream
Member
Emits a value only once quiet passes without a newer one.

Parameters

NameTypeDescription
quietTimeSpan
streamFlowStream<'env, 'error, 'value>

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

keystrokes |> FlowStream.debounce (TimeSpan.FromMilliseconds 300.0)
kind:member

throttle

FlowStream.throttle interval stream
Member
Emits at most one value per interval, keeping the latest value when faster.

Parameters

NameTypeDescription
intervalTimeSpan
streamFlowStream<'env, 'error, 'value>

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

progress |> FlowStream.throttle (TimeSpan.FromMilliseconds 100.0)
kind:member

switchMapFlow

FlowStream.switchMapFlow mapper stream
Member
Maps each value to a flow, interrupting the previous flow when a newer value arrives.

Parameters

NameTypeDescription
mapper'value -> Flow<'env, 'error, 'next>
streamFlowStream<'env, 'error, 'value>

Returns

FlowStream<'env, 'error, 'next>

Verification Examples

queries |> FlowStream.debounce (TimeSpan.FromMilliseconds 200.0) |> FlowStream.switchMapFlow search
kind:member

chunkBySize

FlowStream.chunkBySize size stream
Member
Groups consecutive values into non-empty lists of at most size elements.

Parameters

NameTypeDescription
sizeint
streamFlowStream<'a, 'b, 'c>

Returns

FlowStream<'a, 'b, 'c list>
kind:member

take

FlowStream.take count stream
Member
Emits at most count values.

Parameters

NameTypeDescription
countint
streamFlowStream<'a, 'b, 'c>

Returns

FlowStream<'a, 'b, 'c>

Verification Examples

stream |> FlowStream.take 10
kind:member

skip

FlowStream.skip count stream
Member
Skips the first count values.

Parameters

NameTypeDescription
countint
streamFlowStream<'a, 'b, 'c>

Returns

FlowStream<'a, 'b, 'c>

Verification Examples

stream |> FlowStream.skip 10
kind:member

takeWhile

FlowStream.takeWhile predicate stream
Member
Emits values while a predicate remains true.

Parameters

NameTypeDescription
predicate'a -> bool
streamFlowStream<'b, 'c, 'a>

Returns

FlowStream<'b, 'c, 'a>

Verification Examples

stream |> FlowStream.takeWhile (fun value -> value < 100)
kind:member

skipWhile

FlowStream.skipWhile predicate stream
Member
Skips values while a predicate remains true.

Parameters

NameTypeDescription
predicate'a -> bool
streamFlowStream<'b, 'c, 'a>

Returns

FlowStream<'b, 'c, 'a>

Verification Examples

stream |> FlowStream.skipWhile String.IsNullOrEmpty
kind:member

indexed

FlowStream.indexed stream
Member
Emits each value paired with its zero-based index.

Parameters

NameTypeDescription
streamFlowStream<'a, 'b, 'c>

Returns

FlowStream<'a, 'b, (int * 'c)>

Verification Examples

stream |> FlowStream.indexed
kind:member

scan

FlowStream.scan folder initial stream
Member
Emits successive accumulator states.

Parameters

NameTypeDescription
folder'a -> 'b -> 'a
initial'a
streamFlowStream<'c, 'd, 'b>

Returns

FlowStream<'c, 'd, 'a>

Verification Examples

stream |> FlowStream.scan (+) 0
kind:member

distinctUntilChangedBy

FlowStream.distinctUntilChangedBy projection stream
Member
Suppresses consecutive duplicate values according to a projection.

Parameters

NameTypeDescription
projection'a -> 'b
streamFlowStream<'c, 'd, 'a>

Returns

FlowStream<'c, 'd, 'a>

Verification Examples

stream |> FlowStream.distinctUntilChangedBy id
kind:member

append

FlowStream.append (FlowStream right) (FlowStream left)
Member
Concatenates two streams, evaluating the second only after the first ends.

Parameters

NameTypeDescription
FlowStream rightFlowStream<'a, 'b, 'c>
FlowStream leftFlowStream<'a, 'b, 'c>

Returns

FlowStream<'a, 'b, 'c>

Verification Examples

first |> FlowStream.append second
kind:member

collect

FlowStream.collect mapper stream
Member
Maps each value to a stream and concatenates the resulting streams.

Parameters

NameTypeDescription
mapper'a -> FlowStream<'b, 'c, 'd>
streamFlowStream<'b, 'c, 'a>

Returns

FlowStream<'b, 'c, 'd>

Verification Examples

stream |> FlowStream.collect FlowStream.fromSeq
kind:member

zip

FlowStream.zip rightStream leftStream
Member
Pairs values from two streams until either stream ends.

Parameters

NameTypeDescription
rightStreamFlowStream<'a, 'b, 'c>
leftStreamFlowStream<'a, 'b, 'd>

Returns

FlowStream<'a, 'b, ('d * 'c)>

Verification Examples

left |> FlowStream.zip right
kind:member

runFold

FlowStream.runFold folder initial stream
Member
Folds a stream into one value inside Flow.

Parameters

NameTypeDescription
folder'state -> 'value -> 'state
initial'state
streamFlowStream<'env, 'error, 'value>

Returns

Flow<'env, 'error, 'state>

Verification Examples

stream |> FlowStream.runFold (+) 0
kind:member

runCollect

FlowStream.runCollect stream
Member
Collects all emitted values into a list.

Parameters

NameTypeDescription
streamFlowStream<'a, 'b, 'c>

Returns

Flow<'a, 'b, 'c list>

Verification Examples

stream |> FlowStream.runCollect
kind:member

runDrain

FlowStream.runDrain stream
Member
Consumes a stream and ignores its values.

Parameters

NameTypeDescription
streamFlowStream<'a, 'b, 'c>

Returns

Flow<'a, 'b, unit>

Verification Examples

stream |> FlowStream.runDrain
kind:member

runForEachFlow

FlowStream.runForEachFlow action stream
Member
Runs an effectful action for every stream value.

Parameters

NameTypeDescription
action'value -> Flow<'env, 'error, unit>
streamFlowStream<'env, 'error, 'value>

Returns

Flow<'env, 'error, unit>

Verification Examples

stream |> FlowStream.runForEachFlow save
kind:member

repeatFlow

FlowStream.repeatFlow flow
Member
Creates a stream that emits the result of running flow again for every pull, forever.

Parameters

NameTypeDescription
flowFlow<'env, 'error, 'value>

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

FlowStream.repeatFlow readSensor |> FlowStream.takeWhile (fun reading -> reading.Ok)
kind:member

using

FlowStream.using resource build
Member
Creates a stream that owns a resource: acquired when consumption starts, released when it ends.

Parameters

NameTypeDescription
resourceResource<'env, 'error, 'resource>The resource to acquire.
build'resource -> FlowStream<'env, 'error, 'value>Builds the stream from the acquired resource.

Returns

FlowStream<'env, 'error, 'value>

Verification Examples

let lines path =
    FlowStream.using (Resource.create (Flow.fromBlocking (fun _ -> File.OpenText path)) (fun reader _ -> reader.Dispose(); Task.CompletedTask))
        (fun reader -> FlowStream.repeatFlow (Flow.fromBlocking (fun _ -> reader.ReadLine())) |> FlowStream.takeWhile (isNull >> not))
kind:member

runTryHead

FlowStream.runTryHead stream
Member
Returns the first value, or None for an empty stream, then stops the stream.

Parameters

NameTypeDescription
streamFlowStream<'env, 'error, 'value>

Returns

Flow<'env, 'error, 'value option>

Verification Examples

stream |> FlowStream.runTryHead
kind:member

runTryLast

FlowStream.runTryLast stream
Member
Consumes the whole stream and returns its last value, or None for an empty stream.

Parameters

NameTypeDescription
streamFlowStream<'env, 'error, 'value>

Returns

Flow<'env, 'error, 'value option>

Verification Examples

progress |> FlowStream.runTryLast
kind:member

runCount

FlowStream.runCount stream
Member
Consumes the whole stream and returns how many values it emitted.

Parameters

NameTypeDescription
streamFlowStream<'env, 'error, 'value>

Returns

Flow<'env, 'error, int64>

Verification Examples

stream |> FlowStream.runCount